02 - LangGraph Checkpointer
前置:01 篇的检查点概念。
本篇回答:LangGraph 存了什么、恢复点怎么算出来、崩溃窗口怎么处理、以及它做不到什么。
术语
- 通道(channel ):节点之间传数据的位置,可先理解为状态里的一个字段
- 超级步(superstep):一轮并行执行。执行模型借自 Google Pregel —— 一轮内所有待执行节点并行运行,统一提交后进入下一轮
- thread:一次会话或一个任务的标识,是检查点的主键
一、检查点的数据结构
libs/checkpoint/langgraph/checkpoint/base/__init__.py 中的定义是理解整个机制的入口:
class Checkpoint(TypedDict):
"""某一时刻的状态快照。"""
v: int
# 检查点格式版本号,当前为 1。
# 存在这个字段意味着格式会演进 —— 旧检查点必须能被新代码读懂。
id: str
# 检查点 ID。同时满足唯一且单调递增,因此可直接用于排序。
# 这样列历史时不必依赖时间戳 —— 分布式环境下时钟不可靠
# (节点间漂移、同一毫秒内产生多条)。
ts: str
# ISO 8601 时间戳。仅供展示,不参与排序。
channel_values: dict[str, Any]
# 各通道在该时刻的实际值,即"状态本身"。
channel_versions: ChannelVersions
# 各通道「当前」的版本号。每次写入该通道,版本号递增。
versions_seen: dict[str, ChannelVersions]
# 每个节点"看过"每个通道的哪个版本。
# 与 channel_versions 比对即可算出哪些节点需要执行 —— 见 1.2 节。
updated_channels: list[str] | None
# 本次检查点更新了哪些通道,用于减少下游的比对范围。
1.1 ID 单调递增带来的简化
id 同时承担主键与排序键两个职责。这个决定省掉了一整类问题:列历史记录、找父检查点、判断先后顺序全部退化成整数比较,不需要额外的序号字段,也不受时钟影响。
1.2 恢复点是算出来的,不是存下来的
这是整个机制里最值得理解的一点:LangGraph 没有"执行到第几步"这样的字段。
它存的是版本对照关系:
调度逻辑就是一句话:某节点看过的版本落后于通道当前版本,该节点就进入待执行集合。
这个设计换来三件用"程序计数器"做不到的事:
| 能力 | 为什么成立 |
|---|---|
| 并行分支 | 多个节点各自落后于不同通道,一轮内可全部调度 |
| 人工改状态后自动续跑 | 手工写入某通道 → 版本号递增 → 依赖它的节点自动变为待执行 |
| 图结构变更不影响恢复 | 依赖关系按通道计算,与节点顺序无关 |
可迁移的结论:为 DAG 或状态机做恢复时,记录"每个节点消费到哪个版本"比记录"全局执行到第几步"更健壮 —— 前者是局部信息、天然可并行、对拓扑变化不敏感。
1.3 全量值与增量存储
channel_values 存全量。当某个通道很大(例如一段长对话历史)时,每步存一次会造成显著的写放大 —— 即 01 篇 2.4 节描述的问题。
LangGraph 的解法是 DeltaChannel(当前标注 beta):平时只存增量写入,周期性存一次完整快照作为种子,恢复时从最近快照沿增量链向前叠加。
元数据里记录了触发快照的两个计数:
(updates, supersteps)
updates —— 有多少个超级步写过这个通道
supersteps —— 距离该通道上次快照,总共过了多少个超级步
官方注释说明了第二个计数的作用:
The supersteps bound prevents unbounded ancestor walks on threads where [the channel is rarely written]
很少被写的通道,其"最近快照"可能在几百步之前,恢复时要回溯极长的祖先链。所以除按写入次数触发快照外,还需要一个按超级步数量的强制上限兜底。
这与数据库的 WAL + checkpoint、Redis 的 AOF + RDB 是同一模式:增量存储必须配强制全量点。
二、中间写入:崩溃窗口的处理
01 篇 2.2 节指出,"步骤完成"与"检查点落盘"之间必然存在窗口。LangGraph 用 pending_writes 处理框架内部的这段窗口。
2.1 机制
class CheckpointTuple(NamedTuple):
config: RunnableConfig
checkpoint: Checkpoint
metadata: CheckpointMetadata
parent_config: RunnableConfig | None = None
pending_writes: list[PendingWrite] | None = None
# ↑ 已完成但尚未并入下一个检查点的节点输出
def put_writes(
self,
config: RunnableConfig, # 关联到哪个检查点
writes: Sequence[tuple[str, Any]], # (通道名, 值) 列表
task_id: str, # 哪个节点产生的 —— 恢复时靠它判断该节点是否已完成
task_path: str = "",
) -> None:
"""存储关联到某个检查点的中间写入。"""
方法文档里的关键词是 intermediate:一个超级步内多个节点并行执行,每个节点完成后立即单独落一次写入并带上 task_id;整轮结束后才生成新的 Checkpoint。
没有这个机制,整个超级步的三个节点都要重来。在"一个节点 = 一次模型调用"的场景下,这是三倍的直接成本。
2.2 它的作用范围
pending_writes 保证的是框架不会重复执行已完成的节点。它不能保证节点内部调用外部服务的幂等 —— 如果节点 A 已经发出飞书消息但 put_writes 写入失败,恢复后消息仍会重发。
外部副作用的幂等仍然由 01 篇第三节的幂等键负责,框架层解决不了。
三、接口全貌
BaseCheckpointSaver 的同步方法(每个都有 a 前缀的异步版本):
| 方法 | 作用 | 类别 |
|---|---|---|
get / get_tuple | 读取某个检查点 | 基础 |
list | 列出历史 | 基础 |
put | 写入检查点 | 基础 |
put_writes | 写入中间结果 | 基础 |
delete_thread | 删除某 thread 的全部检查点 | 运维 |
delete_for_runs | 按 run 删除 | 运维 |
copy_thread | 复制整个 thread | 能力扩展 |
prune | 清理 | 运维 |
get_delta_channel_history | 增量通道历史(beta) | 能力扩展 |
前四个是"能存能取",后五个是"能运维、能分支"。两组的存在时间差反映了这类组件的成熟路径:先解决功能,再解决运维。
3.1 copy_thread 支撑的是调试能力
复制整个会话,是"分支"的底层依赖:任务跑到第 18 步,想验证"如果第 10 步换个方向会怎样" —— 复制一份 thread,从第 10 步改。
这是 Agent 调试的高频需求,也是 LangGraph 时间旅行功能的实现基础。
3.2 prune 与 delete_thread 不只是省存储
检查点无限增长。一个 30 步的任务若每步一个检查点、每个检查点携带完整对话历史,存储成本增长很快。
更重要的是:这两个方法同时是数据保留合规的执行点 —— 见 01 篇 2.5 节。上线前不定清理策略,三个月后会同时面对存储账单和合规问询。
四、一致性测试套件
libs/checkpoint-conformance/ 把"什么样的实现算合格的 checkpointer"写成了一套可执行测试。
# libs/checkpoint-conformance/.../validate.py
# 每一项能力对应一套独立测试
_RUNNERS = {
Capability.PUT: run_put_tests,
Capability.PUT_WRITES: run_put_writes_tests,
Capability.GET_TUPLE: run_get_tuple_tests,
Capability.LIST: run_list_tests,
Capability.DELETE_THREAD: run_delete_thread_tests,
Capability.DELETE_FOR_RUNS: run_delete_for_runs_tests,
Capability.COPY_THREAD: run_copy_thread_tests,
Capability.PRUNE: run_prune_tests,
Capability.DELTA_CHANNEL_HISTORY: run_delta_channel_history_tests,
}
async def validate(
registered: RegisteredCheckpointer,
*,
capabilities: set[str] | None = None,
...
) -> CapabilityReport:
"""对一个已注册的 checkpointer 跑验证套件。
capabilities: 指定则只测这些能力;
不指定则「自动探测」实现了哪些,只测实现了的。
"""
这是能力探测式而非全量强制式:实现了哪些就测哪些,未实现的不算失败,产出一份 CapabilityReport 能力表。
4.1 为什么这个设计值得单独学
checkpointer 是典型的"接口简单、语义微妙"组件。这些问题文档很难说清:
list的排序保证是什么?- 并发
put同一个 thread 会发生什么? delete_thread之后get应返回None还是报错?
语义只写在文档里,第三方实现(Redis 版、MongoDB 版、企业自研版)必然各自跑偏,出问题时无法判定是框架缺陷还是实现缺陷。
把规范写成可执行测试,是对齐语义的唯一可靠手段。 这与 Agent 安全 · MCP 授权演进中规范用 MUST / MUST NOT 收紧措辞属于同一类努力。
五、能力边界
| 能力 | 支持情况 |
|---|---|
| 图内部状态恢复 | ✅ |
| 人在环中断与恢复 | ✅ Interrupt 为一等公民,带 ID 可直接恢复 |
| 时间旅行 / 分支调试 | ✅ list + copy_thread |
| 跨进程自动恢复 | ❌ |
| 跨服务长事务 | ❌ 仅覆盖 LangGraph 图内 |
| 定时 / 延迟执行 | ❌ 挂起等待需自行调度 |
| 非 LangGraph 代码 | ❌ |
5.1 最关键的缺口
检查点解决了"状态不丢",没有解决"谁把任务捡起来继续跑"。